Skip to content

feat(connectors): add sink_template and source_template starting points - #3965

Closed
ryerraguntla wants to merge 185 commits into
apache:masterfrom
ryerraguntla:docs/connector-template_code_blog_post
Closed

ryerraguntla wants to merge 185 commits into
apache:masterfrom
ryerraguntla:docs/connector-template_code_blog_post

Conversation

@ryerraguntla

@ryerraguntla ryerraguntla commented Aug 25, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #3956

What

Adds two ready-to-copy connector crates under core/connectors/:

  • core/connectors/sinks/sink_template/
  • core/connectors/sources/source_template/

Each implements every framework-level requirement already expected of a production connector (config parsing with unknown-field rejection, SecretString-wrapped credentials, connectivity validation in open(), retry/backoff + circuit breaker via iggy_connector_sdk::retry, batch processing / Ack-Nack cursor staging, and a canonical test suite), leaving only the backend-specific bits marked TODO(ConnectorDeveloper):

  • sink_template: the outbound data-push call to the destination system.
  • source_template: the connection string/client setup and the data-fetch call against the source system.

Also included:

  • core/connectors/docs/authoring-sinks-and-sources.md — a written guide walking through the required patterns and why each one exists, meant to be cross-linked from the connector READMEs and reused as the basis for an Apache Iggy blog post.
  • Updated .claude/skills/connector-* skill definitions, including a new connector-pr-review skill intended to catch the same recurring review issues (missing retry/circuit-breaker, plaintext secrets, unvalidated config, missing cursor staging) automatically on future connector PRs.
  • Workspace registration (Cargo.toml, Cargo.lock) and README entries (core/connectors/sinks/README.md, core/connectors/sources/README.md,
    core/connectors/README.md) for both new crates.
  • Example runtime configs (core/connectors/runtime/example_config/connectors/{sink,source}_template.toml).

Why

Per the discussion in #3956: connector PRs have been hitting the same review comments repeatedly (improperly typed secrets, missing state-staging, inconsistent error handling, thin test coverage). Rather than repeating that feedback PR after PR, these templates make the required pattern the path of least resistance — a new connector starts from working, reviewed code and only needs the backend-specific logic filled in.

Local Execution
Passed
Pre-commit hooks ran

Testing

  • cargo fmt --all
  • cargo clippy --all-targets --all-features -- -D warnings
  • cargo build -p iggy_connector_template_sink -p iggy_connector_template_source
  • cargo test -p iggy_connector_template_sink -p iggy_connector_template_source
  • cargo machete
  • cargo sort --workspace
  • typos
  • ./scripts/ci/license-headers.sh --check

AI Usage
Claude
Implementation was AI-generated with human review
Each line of the code is reviewed and can be explained by developer

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Aug 25, 2026
@ryerraguntla
ryerraguntla requested review from hubcio and a lite review from Copilot August 25, 2026 15:46

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Warning

Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.

Adds ready-to-copy “template” sink and source connector crates (plus docs and skill checklists) to standardize connector implementations and reduce recurring review feedback (secrets handling, retries/circuit-breaker, config validation, cursor staging, and canonical tests).

Changes:

  • Introduces sink_template and source_template crates implementing the expected connector scaffolding with TODO(ConnectorDeveloper) placeholders for backend-specific logic.
  • Adds example runtime configs and README/docs entries to make the templates discoverable and easy to adopt.
  • Updates .claude/skills guidance (including new connector-pr-review) to encode recurring review blockers and conventions.

Reviewed changes

Copilot reviewed 21 out of 22 changed files in this pull request and generated 7 comments.

Show a summary per file
File Description
core/connectors/sources/source_template/src/lib.rs New source template implementation + cursor staging + tests
core/connectors/sources/source_template/config.toml Example source template connector config
core/connectors/sources/source_template/README.md Source template usage guidance and expectations
core/connectors/sources/source_template/Cargo.toml Source template crate metadata/deps for workspace + cdylib
core/connectors/sources/README.md Registers source_template in sources index
core/connectors/sinks/sink_template/src/lib.rs New sink template implementation + batching/circuit breaker + tests
core/connectors/sinks/sink_template/config.toml Example sink template connector config
core/connectors/sinks/sink_template/README.md Sink template usage guidance and expectations
core/connectors/sinks/sink_template/Cargo.toml Sink template crate metadata/deps for workspace + cdylib
core/connectors/sinks/README.md Registers sink_template in sinks index
core/connectors/runtime/example_config/connectors/source_template.toml Runtime example config for source template
core/connectors/runtime/example_config/connectors/sink_template.toml Runtime example config for sink template
core/connectors/BLOG_POST.md Draft blog post announcing the templates
Cargo.toml Adds both template crates to workspace members
.claude/skills/connectors-overview/SKILL.md Updates connector overview skill (incl. review checklist pointers)
.claude/skills/connector-source/TEMPLATE.md Reworks source template kit content to align with ACK/NACK + patterns
.claude/skills/connector-source/SKILL.md Documents ACK/NACK contract and staging rules more explicitly
.claude/skills/connector-sink/TEMPLATE.md Reworks sink template kit content and pre-flight guidance
.claude/skills/connector-sink/SKILL.md Adds pointers to PR pre-flight skill
.claude/skills/connector-runtime/SKILL.md Updates metrics naming guidance
.claude/skills/connector-pr-review/SKILL.md New PR review checklist skill for connector PRs

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread core/connectors/sources/source_template/src/lib.rs
Comment thread core/connectors/sinks/sink_template/src/lib.rs
Comment on lines +300 to +301
self.config.retry_max_delay.as_deref(),
DEFAULT_RETRY_MAX_DELAY,
// and unknown keys are rejected outright so a typo in a TOML file fails at
// load time instead of silently doing nothing.

#[derive(Debug, Clone, Serialize, Deserialize)]
Comment on lines +92 to +93
#[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")]
pub connection_string: SecretString,
Comment on lines +100 to +104
#[serde(
default,
serialize_with = "iggy_common::serde_secret::serialize_optional_secret"
)]
pub auth_token: Option<SecretString>,
Comment on lines +571 to +581
#[tokio::test]
async fn poll_returns_empty_without_error_when_circuit_is_open() {
let source = TemplateSource::new(1, test_config(), None);
source.circuit_breaker.record_failure().await;
source.circuit_breaker.record_failure().await;
source.circuit_breaker.record_failure().await;
assert!(source.circuit_breaker.is_open().await);

let result = source
.poll()
.await
@codecov

codecov Bot commented Aug 25, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 84.97%. Comparing base (996ac04) to head (ec3b463).
⚠️ Report is 162 commits behind head on master.

Additional details and impacted files
@@             Coverage Diff              @@
##             master    #3965      +/-   ##
============================================
+ Coverage     84.76%   84.97%   +0.21%     
- Complexity     1405     1575     +170     
============================================
  Files          1224     1236      +12     
  Lines        177971   180844    +2873     
  Branches     144285   144501     +216     
============================================
+ Hits         150854   153671    +2817     
+ Misses        23100    23097       -3     
- Partials       4017     4076      +59     
Components Coverage Δ
Rust Core 85.67% <ø> (+0.05%) ⬆️
Java SDK 68.68% <ø> (+1.33%) ⬆️
C# SDK 77.41% <ø> (+2.02%) ⬆️
Python SDK 90.97% <ø> (+0.90%) ⬆️
PHP SDK 85.67% <ø> (+0.02%) ⬆️
Node SDK 96.36% <ø> (+0.23%) ⬆️
Go SDK 70.20% <ø> (+0.94%) ⬆️
Files with missing lines Coverage Δ
core/ai/mcp/src/log.rs 14.42% <ø> (+1.08%) ⬆️
core/ai/mcp/src/main.rs 56.57% <ø> (+3.04%) ⬆️
core/ai/mcp/src/service/mod.rs 62.71% <ø> (+1.13%) ⬆️
core/ai/mcp/src/service/requests.rs 100.00% <ø> (ø)
core/binary_protocol/src/codes.rs 100.00% <ø> (ø)
core/binary_protocol/src/consensus/header.rs 81.87% <ø> (ø)
core/binary_protocol/src/consensus/operation.rs 97.54% <ø> (ø)
core/binary_protocol/src/consensus/reply_result.rs 100.00% <ø> (ø)
core/binary_protocol/src/dispatch.rs 92.70% <ø> (ø)
core/binary_protocol/src/primitives/ack_level.rs 97.61% <ø> (ø)
... and 100 more

... and 141 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

ryerraguntla and others added 22 commits August 25, 2026 13:10
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
….com/ryerraguntla/iggy into docs/connector-template_code_blog_post
…apache#3970)

Every journal flush issued two serialized fdatasyncs, log then
index, so an fsync-gated topic paid two device round trips per
ack. Both files now persist under one futures::future::join:
13.4k to 26.1k msg/s and p50 5.83 to 3.01 ms on ext4 with
enforce_fsync and messages_required_to_save = 1. Cursors advance
only after both saves succeed, so a failed half keeps its slot
and a later flush rewrites the same positions instead of
appending a duplicate.

Removing the write barrier forced a new recovery contract. An
index entry is derived from the log and, without enforce_fsync,
proves nothing about it, so recovery checksum-walks every
segment from byte 0 and keeps a clean index untouched. Every
entry must match its log batch's offset, timestamp, partition
and extent, or the index is rebuilt: ascending columns alone let
a corrupt entry point at the wrong valid batch and polls
silently skipped offsets. Under enforce_fsync serialized flushes
prove the prefix below the last entry, so a healthy boot still
anchors there; a gap deeper than the one in-flight chunk refuses
as FsyncedLogLoss, and a rebuild stopping below the last entry's
position refuses as FsyncedRebuildShortfall rather than
truncating bytes a completed flush made durable. Probe-budget
exhaustion always refuses, since recovering as empty would serve
the segment and re-mint offsets over bytes never scanned.
Previously any index the log outran refused boot outright and
tombstoned single-replica partitions with healthy logs.

A persist failure on a cluster-committed op panicked on the
shard pump, where compio::runtime::spawn swallows the unwind:
the process stayed up with every partition on that shard dead.
The partition now fences itself, the pump breaks, flips the
shared shutdown flag before its final flush (the flush writes to
the failed device and can stall forever; the flag is what arms
the watchdog and the bounded drain), and a failed shutdown flush
fences its partition too, so the exit is ShardFatal and non-zero
rather than a clean report over unpersisted committed data.

Also: the active segment fills the same read-state slot sealed
segments use, ending one openat per head poll; latency extremes
backfill only as a whole group; segment writer save methods are
crate-private; index scans yield on a stride.
Merging adjacent producer batches released their byte permits before
the write completed, allowing buffered and in-flight data to exceed
max_buffer_size. Large permit totals could also overflow during the
merge and permanently reduce available capacity.

Retain each permit separately until send_internal completes. Validate
batch sizes before u32 conversion, split merged writes at the wire
limit, and handle zero buffer and in-flight limits without invalid
semaphores.

Update max_buffer_size documentation and add regression coverage for
permit lifetime, unlimited limits, and budgets above u32::MAX. Keep
the edge.6 package versions synchronized.

Fixes apache#3932
…he#3958)

Our partition operations did not use the `request` counter which meant
that partition ops did not participate in the deduplication process.
This PR makes the partition operations bump the request counter for all
of the SDKs
…up (apache#4007)

A partition group with no history of its own seeds its view and
log_view from the metadata view this replica happens to publish.
Nothing bounds that view by the group's real one, so a replica that
materialises the group late, after a further metadata election,
starts Normal on an empty log above every peer. That empty log is
canonical in the next DVC merge: as a backup it collects a nack
quorum against committed ops and the group deadlocks for good, the
superblock persisting the view so a restart re-wedges; as the
primary-by-index it answers the real primary's heartbeat with a
StartView at op 0 that resets the peers' heads under their commit
floor, and every new write is fenced. Introduced with the seed in
apache#3985.

Record the view the admitting primary mints the create in on the
committed metadata entry and seed from that instead. Every replica
reads the same value, and the group's view never drops below it, so
an empty log at the creation view is the floor a view-0 start used
to be, never canonical over committed history. The reconciler no
longer needs the published-view sentinel, its deferral, or the
double atomic read that came with it.

The value rides the request body, not the prepare header: a
header's view is the one field a post-view-change retransmit
restamps (replicate_preflight refuses the frame otherwise), so a
header-derived value is a per-replica property of delivery history
and durably diverges the metadata STM. The body is sealed by
checksum_body and never restamped; entries journaled before the
trailing field decode it as 0, the safe floor.

Chosen over electing into the seed: no extra round trip per group,
no consensus change, and a late materialiser joins as a backup at
or below its peers instead of forcing an election it may win with
an empty log. The snapshot format moves to version 5 with the new
trailing field defaulted, so version 3 and 4 checkpoints still
decode.
@hubcio

hubcio commented Sep 21, 2026

Copy link
Copy Markdown
Contributor

moving this to draft to not trigger stale bot

@hubcio
hubcio marked this pull request as draft September 21, 2026 08:56
@github-actions github-actions Bot removed the S-waiting-on-review PR is waiting on a reviewer label Sep 21, 2026
hubcio and others added 24 commits September 21, 2026 11:31
…#4224)

The link shard that owns a peer's replica socket burns a full
core under sustained replicated writes. Nothing batches its reads:
read_message issues one io_uring read for the 256-byte header and
a second for the body. Bytes already sitting in the kernel receive
queue wait for the next call.

A read-ahead buffer sized by message_bus.replica_read_buffer_size
lets one socket read serve a whole burst. It adds no latency, because
the reader never waits for the buffer to fill. One fill delivers
at most one buffer, so a body above that size crosses the buffer in
buffer-sized pieces and pays one extra pass over the payload and one
read per piece. A read that large goes straight into the frame's own
buffer instead.

Zero keeps the unbuffered path, so the A/B baseline is a config change
and not a separate build. The plaintext writer already coalesces into
one writev, so it is left alone.

replica_socket_reads_total and replica_inbound_frames_total carry the
evidence on the real workload: their ratio is the batching factor,
and only link shards bump them. A read counts once it completes,
so a link torn down mid-read cannot inflate the ratio.
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
# Conflicts:
#	.claude/skills/connector-source/TEMPLATE.md
apache#4104 landed on master after this branch's commits were written and
replaced the connectivity retry API: ConnectivityConfig is gone,
check_connectivity_with_retry now takes an owned RetryPolicy. Map the
old fields onto the new struct with identical semantics
(max_open_retries -> max_attempts, retry_delay -> base_delay,
open_retry_max_delay -> max_delay) in both templates.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JpGKmaGSG6RdBUiaq3BW2f
….com/ryerraguntla/iggy into docs/connector-template_code_blog_post
@hubcio

hubcio commented Sep 23, 2026

Copy link
Copy Markdown
Contributor

@ryerraguntla looks like merge is broken in this PR
image

@hubcio

hubcio commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

history broken in this PR, please reopen a new one as fresh

@hubcio hubcio closed this Sep 28, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Add sink/source connector templates under core/connectors/ as fill-in-the-blank starting points